iT邦幫忙

2026 iThome 鐵人賽

DAY 15
0

不知不覺鐵人賽走到過半了,標題那句「1+1+1>3 ~ Spark 與 DataFusion、Comet 效能煉金術 ~」裡的第三個 1,那顆叫 Comet 的行星,已經躺在標題上 14 天了,動也不動,今天要來介紹什麼是 Apache DataFusion Comet?

那回顧一下前面 14 天介紹了什麼?
D1-D5 拆硬體跟 Arrow,D6-D8 過場(codegen 對照、datafusion-cli 起手、SortExec RowFormat 實驗),D9-D12 講儲存格式(Parquet、Iceberg、晚物化、delete files),D13 問可嵌入引擎為什麼還要再做一個,D14 拉了 morsel-driven 論文,而這些全都是 Comet 的前情提要,前提舖完了,主角該登場。

我目前是 Apache DataFusion Comet 開源專案的貢獻者,今天先畫一張地圖:先一句話說 Comet 是什麼、它要解決什麼痛點,再對照 Photon / Gluten / Comet 三條路,最後標出六個接點,後面幾天會一個個回來理解。

另外這邊的版本固定用 Comet 1.0.0 / DataFusion 54.1 / Arrow 58.4
因為 Comet 變化太快,所以先說明一下

一句話 Comet 是什麼

先來個一句話定義:

Comet 是 Apache Spark 的執行後端(execution backend):Spark 的前端、Catalyst 分析、優化、AQE 全部留著,只把 physical plan 換成 DataFusion 的 ExecutionPlan,在同一個 executor JVM 裡跑 native code

拆三個關鍵詞:

  • 執行後端,不是 Spark 替代品
    Spark app、SparkSession、DataFrame API 完全不動,Comet 是插件(plugin),註冊一個叫 CometSparkSessionExtensions,擋在 physical planning 之後、實際執行之前

  • DataFusion,不是自己從頭寫
    D13 提到可嵌入引擎的動機就是為這件事鋪路,Comet 把 hash join、sort、vectorized filter 這些通用運算交給 DataFusion,但也不是只有 glue code,像是 JSON / regex / URL parsing、date/time 與 decimal kernel、string 與 array kernel 這些 Spark 語義獨有的部分,Comet 也有自己實作

  • 在同一個 executor JVM 裡,不是外接 process
    Spark task 在 JVM 裡跑起來、Comet 把資料指標交給 Rust、Rust 算完把結果指標交回 JVM,中間走 Arrow C Data Interface 零拷貝(Zero Copy),走 JNI 邊界,這也是為什麼記憶體管理會變得很棘手

Comet 要解決什麼

寫這系列碰到最多的一個問題就是「Spark 已經很成熟了,為什麼還要 Comet?」講三個實際痛點。

第一個痛點:JVM 稅
Spark 從 2.0 之後靠 whole-stage codegen(WSCG)已經很努力地把運算子 JIT 掉,但 JIT + GC + JVM 記憶體佈局,跟 native vectorized 執行比還是差一截。Photon SIGMOD 2022 的 microbenchmark 給了個大致範圍,本季不做效能斷言就不引具體倍數,但方向是明確的:熱點運算子(scan / filter / projection / hash aggregation)native 可以贏,可是重寫成本很高

第二個痛點:重寫成本沒人養得起
Databricks 花了幾年、幾十個 SIGMOD 級工程師從頭寫 C++ 的 Photon,證明了 native 可以贏,但也證明了這條路開源社群跟不上。想從零寫一個可以取代 Spark 執行引擎的 C++ 系統,得同時有一群 Spark 核心貢獻者跟一群 C++ 高手,還要能維持好幾年不做別的事。

第三個痛點:生態鎖定風險
Photon 只在 Databricks Runtime 跑,想用它得買 Databricks 的產品,開源社群想要 native 執行的收益,可是不想付只能在某家雲上跑的代價。

Comet 的取捨是:用 DataFusion 當現成的執行引擎,用 protobuf IR + Arrow C Data Interface + JNI 三層膠水,把它嵌回 Spark。這就是可嵌入引擎的動機在工業界的落地,DataFusion 不是要贏 Spark,是要當 Spark 底下那塊可以被替換的執行引擎。

前世今生:三條路的對照

Comet 不是第一個想 把 native 執行接回 Spark 的專案,前面走過兩條路,跟 Comet 湊起來剛好是三個對照組。

Photon(Databricks,SIGMOD 2022):C++ 從頭寫,跟 Databricks Runtime 深度整合,優點是極致效能、跟 Spark SQL 語義完全對得起來,代價是只在 Databricks 跑、閉源、養團隊的成本嚇死人,它的貢獻是證明 native 執行可以贏,但無法開源複製

Gluten(2021 年由 Intel 主導,2026 年 3 月畢業為 Apache 專案):走 Substrait 當跨語言 IR,後端可以接 Velox(Meta 的可嵌入 C++ 引擎)或 ClickHouse,優點是後端可換、有 Substrait 標準加持,但代價是 Substrait 核心規格外要靠 simple extensions 機制去補 Spark 的特殊語義(某些型別轉換、NULL 處理的 corner case),extension 一多、跨引擎相容的承諾就會被稀釋;而且跨後端切換的複雜度會落到使用者頭上——你到底用 Velox 還是 ClickHouse,選錯了就是重跑一次。它的位置是走「標準 IR」路線,追求跨引擎相容

Comet(Apache DataFusion Comet,2024 年 3 月捐給 Apache DataFusion 成為它的 subproject,五個月後發 0.1.0;2026 年 8 月 1.0.0):自訂 protobuf 當 Spark ↔ DataFusion 的 IR,不用 Substrait;後端固定綁 DataFusion,Rust 寫的,優點是 protobuf schema 完全長成 Spark plan 的形狀,Spark 特殊語義能一對一對應;DataFusion 是 Rust,開發速度比 C++ 快,社群長得比 Velox 快,但代價是 protobuf schema 是 Comet 內部契約,不能跨引擎重用;綁在 DataFusion 上,DataFusion 沒實作的表達式就要 fallback 回 JVM。它的位置是放棄跨引擎相容,換取 Spark 語義的一比一對應跟開發效率

那三個核心架構差在哪呢?
Photon 是走一條 Catalyst rule,從 file scan 節點由下往上走 plan,把支援的節點換成 Photon 節點;遇到不支援的節點就插入 transition 節點,把 columnar 資料轉回 row-wise 給 legacy Spark 跑,而這跟 Comet / Gluten 的攔截點是同一層。

        Spark User Code
              │
        Catalyst Analyzer
              │
        Catalyst Optimizer
              │
        Physical Planner
              │
    ┌─────────┴─────────────────────────────────┐
    │                                           │
    │  Comet plugin ─────────  攔在 physical plan:protobuf IR → DataFusion (Rust)
    │                                           │      → vanilla Spark 佈署都能裝
    │                                           │
    │  Gluten plugin ────────  攔在 physical plan:Substrait IR → Velox / ClickHouse (C++)
    │                                           │      → vanilla Spark 佈署都能裝
    │                                           │
    │  Photon (DBR only) ────  攔在 physical plan:Catalyst rule 由下往上替換 → C++ 引擎
    │                                           │      → 深度嵌 Databricks Runtime、只在 DBR 跑
    │                                           │
    └─────────┬─────────────────────────────────┘
              │
        Physical Execution
              │
        Task Scheduler & Shuffle
              │
        RDD & Storage

真正的差異在整合深度,不在替換邊界:Photon 是嵌進 Databricks Runtime 的 C++ 引擎(每個 pipeline instance 在 task 內單執行緒跑,平行度靠 Spark 的 task model;Comet 在這一層也是一樣,一個 Spark task 一個執行實例),透過 JNI 呼叫、跟 DBR 共用同一套查詢優化器與 adaptive optimization,並且跟 DBR 的規劃器與記憶體管理器綁得很深。Comet 跟 Gluten 則是 vanilla Spark 上的 plugin,任何 Spark 佈署都能裝;Photon 只在 DBR 跑,差別是有沒有 fork runtime、整合到什麼程度

這也是為什麼今天標題副標是那顆躺了 14 天的行星該登場了:D13 的可嵌入引擎、D14 的 morsel-driven,兩天都是在給 Comet 的替換邊界做鋪墊,Spark 的分析、優化、AQE 全部留著,physical plan 之後才換 DataFusion

六個接點:整段 Comet 的地圖

Comet 是黏在 Spark 跟 DataFusion 中間的一層膠水,上面那層是 JVM、Scala、列狀資料;下面那層是 Rust、native、Arrow 欄狀資料,兩邊沒有一樣東西是天生相容的,所以每一樣要跨過去的東西都得有個交接的地方,計畫、資料、記憶體、責任邊界,一共六個

六個放在一起看,其實就是一次完整的查詢生命週期:計畫進來、翻譯過去、資料出來、記憶體收拾掉;中間遇到接不住的就退回去,跨 stage 就走 shuffle。

接點一:plan 攔截——交接的是執行計畫的控制權

Spark 把查詢排成一棵執行計畫樹之後,Comet 要有機會在它真的跑起來之前先攔一手。做法是把兩條規則掛進 Spark 的規劃流程,走過樹上每個節點問一句:這個我能接嗎?能接就換成 Comet 自己的版本,接不了就原封不動留給 Spark。

掛的位置是 CometSparkSessionExtensions,用的是 injectColumnarinjectQueryStagePrepRule 這兩個 Spark 的擴充點;會被換掉的典型節點是 scan、filter、project、aggregate。

順帶解掉一個常見疑惑:AQE 會在執行到一半重新規劃,那 Comet 怎麼辦?答案是不用怎麼辦,規則掛在 query stage prep 上,AQE 每重新規劃一次,這條規則就自動再跑一次。

接點二:protobuf IR——交接的是計畫本身

計畫樹在 JVM 裡,真正要執行它的引擎在 Rust 裡,要把一棵樹從 JVM 送到 Rust,得先變成兩邊都看得懂的格式,Comet 用的是 protobuf。

接點三:Arrow C Data Interface——交接的是資料

Rust 算完一批資料,要交回 JVM,資料量很大複製一份的成本不能接受,所以要交的只是資料放在記憶體哪裡這個資訊。

Arrow C Data Interface 就是這件事的約定:兩個 C 結構體(ArrowArrayArrowSchema)描述資料的位置和長相,再附一個 release callback,讓拿到資料的一方用完之後回報一聲,所有權就靠這個 callback 交接。

接點四:JNI 邊界——交接的是記憶體的所有權跟生命週期

callback 講清楚了誰該收拾,但沒講清楚什麼時候。三件麻煩事:

Rust 這邊一批資料什麼時候算用完,不是 JVM 說了算,所以 release 到底何時被觸發並不直觀。JVM 的 off-heap 記憶體池跟 Rust 的記憶體池是兩本帳,得想辦法合起來算,否則兩邊都以為自己還有預算,最後一起爆掉。還有 Rust 這邊是 async、JVM 那邊是 blocking,兩種模型要怎麼交手。

留一個現成的線索:Comet 0.17 在 CometPlainVector 裡把變寬向量的 offsetBufferAddress 快取起來,不再每次存取都重算,會特地做這個優化,代表跨 JNI 要一次位址並不便宜。

前四個接點講的是資料與控制權怎麼過界,接下來的兩個不太一樣:一個是過不去的時候怎麼退,一個是怎麼插進既有的基礎設施,但兩邊仍然需要做協調,所以還是接點。

接點五:Fallback 粒度——交接的是執行責任的邊界

遇到一個 Comet 算不了的表達式(DataFusion 沒實作,或是 Rust 版本在邊界情況跟 Spark 對不起來),要退多少回去給 Spark?

以前是整塊退,一個表達式不會算,它所在的那整棵子樹都退回 Spark 跑,資料還得從 Arrow 的欄狀格式轉成 Spark 的列狀格式再轉回來,來回一趟的成本常常把前面省下的都吃掉。

0.17 引入 codegen dispatcher 換了個做法:資料留在 Comet 的管線裡不動,只針對那一個表達式呼叫 Spark 自己產生的程式碼,直接作用在 Arrow batch 上。Scala 跟 Java UDF 也走這條路。同一版新增支援的 Spark 表達式一口氣到 120+ 個,其中一部分走 dispatcher、一部分是新寫的 Rust 實作。

1.0 再把 cast 這個大戶從整塊退搬到 dispatch,包含 spark.sql.legacy.castComplexTypesToString.enabled 這類 legacy 設定的變體。現在 extended explain 會直接告訴你,這個查詢裡哪些走 native、哪些走 dispatch。

先做通用機制、再把過去只能退回去的大戶搬過來

接點六:Shuffle——交接的是跨節點搬資料的基礎設施

前四個接點是執行時臨時決定要不要走 native;Shuffle 不一樣
Comet 要插進 Spark 的 shuffle 機制,得註冊自己的 CometShuffleManager,而 spark.shuffle.manager 是 Spark 的靜態設定,必須在 SparkContext 建起來之前就設好,開跑之後改不了。這道縫的決定被迫推到連 context 都還沒起來的時刻,是整段 Comet 裡替換邊界最前面的一處。

Shuffle 的退場機制是三段式的,剛好是接點五在運算子層的鏡像:先試 Native Shuffle,它的分區鍵只能是 primitive 型別;不行就改試 Columnar Shuffle,這個雖然是 JVM 實作,但仍然是 Comet 的,分區鍵可以是複雜型別;兩個都不適用,才真的把 shuffle 還給 Spark。

Comet 1.0.0 現在能做什麼

給個簡短的現況清單:

  • Spark 支援(以 1.0.0 為準):Spark 3.4 到 4.1 支援、4.2 實驗性;完整支援 Spark 4.0 起預設開啟的 ANSI mode;CI 跑 3.4.3 / 3.5.9 / 4.0.4 / 4.1.3。JDK 11 與 Spark 3.4 在 1.0.0 已標 deprecated,1.1.0 移除——想安穩留在 1.x 就趁早升
  • Iceberg(以 1.0.0 為準):跟 Iceberg 1.11 整合,V1 / V2 是目前的主流支援路徑;V3 目前只有第一個特性全表加密可用,deletion vector 與 VARIANT 等其他 V3 特性優雅降級——D20 會回扣「優雅降級在規格層怎麼保證」
  • Shuffle:native + JVM 兩個模式,可以按 workload 切
  • Fallback:0.17 引入 codegen dispatcher(Scala / Java UDF 走這條路;同版一次新增 120+ 個 Spark 表達式支援,其中一部分走 dispatcher、一部分是新的 Rust 實作),1.0 再把 cast(含 legacy config 變體)從 fallback 搬到 dispatch;extended explain 會標示 native / dispatch 的覆蓋率
  • 平台:Maven Central 上的 jar 只 bundle Linux 的 native library(amd64 / arm64),macOS 使用者要自己 build from source。amd64 build 用 x86-64-v3 target(含 AVX2)、arm64 用 neoverse-n1,太舊的 CPU 跑會 SIGILL——寫效能主題的人自己踩到過就會記得這件事

不能做的(或做不好的)也講清楚:某些 Spark 特殊型別(TimestampNTZ 的部分行為、Decimal 高精度)、部分 window function、一些 catalyst-only 的 UDF 還是要 fallback。這些正是「六個接點」裡的成本面。

怎麼跳進來貢獻

Comet 是 Apache DataFusion 底下的 subproject,給幾個具體入口:

可以貢獻的類型從輕到重:

  1. 跑 benchmark 回報 fallback 熱點
    拿 TPC-H / TPC-DS / ClickBench 跑一次,看哪些查詢 fallback 到 Spark、哪些節點是熱點,回報到 issue tracker 就是很有價值的訊號
  2. 寫測試蓋 edge case
    NULL 處理、型別轉換、時區這幾塊是重災區,找一個 Spark 有 Comet 還不完全一致的地方,寫測試把它釘死
  3. 補一個 SQL 表達式
    看 Comet 目前 fallback 哪些表達式、DataFusion 有沒有對應實作,如果 DataFusion 有就補 protobuf serialization + Rust dispatch;如果 DataFusion 沒有就要先在 DataFusion 那邊補
  4. 文件 / release note

那就明天見~

D15 到這裡告一段落,Comet 這顆行星終於浮出來啦!
之後每一天都會回推文章在comet中發揮了什麼意義

明天 D16 接執行模型,兩個題目壓在一起講:Volcano pull-based 執行為什麼在 DataFusion 上還活得好好的,以及 Tokio 這個本來是給網路 I/O 用的 async runtime 怎麼被拿來排 CPU 密集的查詢工作。
回扣點會落在接點四 JNI 邊界
因為 Rust 這邊是 async、JVM 那邊是 blocking,兩者怎麼交手是 Comet 的第一個工程難題。

留一個沒有標準答案的思考題:Comet 選了 Java 規劃 + Rust 執行 這條路,為什麼不是全部沿用(原本 Spark,慢但穩)也不是全部重寫(Photon,快但重)?替換邊界在 physical plan 之後,表示哪些決策要留在 Spark、哪些搬到 Rust?
(提示:AQE 的動態重規劃、cost-based optimizer 的統計資訊、UDF 的 codegen、shuffle 的分區策略——這幾個要住在哪一邊?)

那就明天見~

參考資料


上一篇
Day 14 Morsel-Driven Parallelism:DataFusion 多核並行與 work-stealing
下一篇
Day 16 DataFusion 執行模型:Volcano Pull、RecordBatch 與 Tokio Async
系列文
1+1+1>3 ~ Spark 與 DataFusion、Comet 效能煉金術 ~21
圖片
  熱門推薦
圖片
{{ item.channelVendor }} | {{ item.webinarstarted }} |
{{ formatDate(item.duration) }}
直播中

尚未有邦友留言

立即登入留言